package com.flink.com.source

import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment


object HdfsFileSource {
  def main(args: Array[String]): Unit = {
    val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment
    env.setParallelism(3);
    //2、导入隐式转换
    import org.apache.flink.streaming.api.scala._
    val streamSource: DataStream[String] = env.readTextFile("hdfs://node01.com:8020/wc.txt")
    val result = streamSource.flatMap(_.split(" ")).map((_, 1)).keyBy(0).sum(1)
    result.print()
    env.execute("du qu hdfs")

  }

}
